package cn.dataling.dag.task;

import cn.dataling.dag.enums.DagWorkflowNodeState;
import cn.dataling.dag.pojo.DagWorkflowNode;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.util.concurrent.TimeUnit;

public class FlinkDagWorkflowTask implements DagWorkflowNodeTask {

    private static final Logger LOGGER = LoggerFactory.getLogger(FlinkDagWorkflowTask.class);

    @Override
    public void run(DagWorkflowNode dagWorkflowNode) {
        try {
            TimeUnit.SECONDS.sleep(5);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        LOGGER.info("{} 任务正在执行", dagWorkflowNode.getId());

        dagWorkflowNode.setState(DagWorkflowNodeState.COMPLETED);
    }
}
